🎖️GitЯра🎖️
Node / meshtastic / Meshtastic-Android / files / core / data / src / commonMain / kotlin / org / meshtastic / core / data / manager / PacketHandlerImpl.kt
Displaying Raw • Download
core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/PacketHandlerImpl.kt bc4e9da3ad7a2a22768cbaecc75810e018bc40e2 (bc4e9da3) Text, 16.71 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.data.manager
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.coroutines.CancellationException
Tff7b72import T7ee787kotlinx.coroutines.CompletableDeferred
Tff7b72import T7ee787kotlinx.coroutines.Deferred
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.TimeoutCancellationException
Tff7b72import T7ee787kotlinx.coroutines.asDeferred
Tff7b72import T7ee787kotlinx.coroutines.delay
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787kotlinx.coroutines.withTimeout
Tff7b72import T7ee787kotlinx.coroutines.withTimeoutOrNull
Tff7b72import T7ee787org.koin.core.annotation.Single
Tff7b72import T7ee787org.meshtastic.core.common.di.ServiceScope
Tff7b72import T7ee787org.meshtastic.core.common.util.handledLaunch
Tff7b72import T7ee787org.meshtastic.core.common.util.nowMillis
Tff7b72import T7ee787org.meshtastic.core.model.ConnectionState
Tff7b72import T7ee787org.meshtastic.core.model.DataPacket
Tff7b72import T7ee787org.meshtastic.core.model.MeshLog
Tff7b72import T7ee787org.meshtastic.core.model.MessageStatus
Tff7b72import T7ee787org.meshtastic.core.model.RadioNotConnectedException
Tff7b72import T7ee787org.meshtastic.core.model.util.toOneLineString
Tff7b72import T7ee787org.meshtastic.core.model.util.toPIIString
Tff7b72import T7ee787org.meshtastic.core.repository.ConnectionStateProvider
Tff7b72import T7ee787org.meshtastic.core.repository.MeshLogRepository
Tff7b72import T7ee787org.meshtastic.core.repository.PacketHandler
Tff7b72import T7ee787org.meshtastic.core.repository.PacketRepository
Tff7b72import T7ee787org.meshtastic.core.repository.RadioInterfaceService
Tff7b72import T7ee787org.meshtastic.proto.FromRadio
Tff7b72import T7ee787org.meshtastic.proto.MeshPacket
Tff7b72import T7ee787org.meshtastic.proto.QueueStatus
Tff7b72import T7ee787org.meshtastic.proto.Routing
Tff7b72import T7ee787org.meshtastic.proto.ToRadio
Tff7b72import T7ee787kotlin.time.Duration
Tff7b72import T7ee787kotlin.time.Duration.Companion.milliseconds
Tff7b72import T7ee787kotlin.time.Duration.Companion.minutes
Tff7b72import T7ee787kotlin.time.Duration.Companion.seconds
Tff7b72import T7ee787kotlin.uuid.Uuid
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooManyFunctionsTa5d6ff"Tb4b4b4)
Tf0883e@Single
Tff7b72class T56d364PacketHandlerImplTb4b4b4(
Tff7b72private Tff7b72val Te6edf3packetRepositoryTb4b4b4: Te6edf3LazyTff7b72<Te6edf3PacketRepositoryTff7b72>Tb4b4b4,
Tff7b72private Tff7b72val Te6edf3radioInterfaceServiceTb4b4b4: Te6edf3RadioInterfaceServiceTb4b4b4,
Tff7b72private Tff7b72val Te6edf3meshLogRepositoryTb4b4b4: Te6edf3LazyTff7b72<Te6edf3MeshLogRepositoryTff7b72>Tb4b4b4,
Tff7b72private Tff7b72val Te6edf3connectionStateProviderTb4b4b4: Te6edf3ConnectionStateProviderTb4b4b4,
Tff7b72private Tff7b72val Te6edf3scopeTb4b4b4: Te6edf3ServiceScopeTb4b4b4,
Tb4b4b4) Tb4b4b4: Te6edf3PacketHandler Tb4b4b4{
Tff7b72companion Tff7b72object Tb4b4b4{
Tff7b72private Tff7b72val Te6edf3TIMEOUT Tff7b72= T79c0ff5.Te6edf3seconds
T8b949e/**
* Grace period after which a sent packet still [MessageStatus.ENROUTE] is stamped [Routing.Error.TIMEOUT]
* (retryable) instead of showing as sending forever. Generous — well past the radio's retransmit window — and
* matches iOS's sendAckTimeout so both apps time out alike.
*/
Tff7b72internal Tff7b72val Te6edf3SEND_ACK_TIMEOUT Tff7b72= T79c0ff5.Te6edf3minutes
T8b949e/**
* Minimum re-arm delay on reconnect: the firmware's phone-queue backlog may still deliver the missing ACK/NAK
* just after the config handshake, so give it a moment before stamping a timeout.
*/
Tff7b72internal Tff7b72val Te6edf3REARM_GRACE Tff7b72= T79c0ff3T79c0ff0.Te6edf3seconds
T8b949e/**
* Firmware-internal `ErrorCode` (MeshTypes.h `ERRNO_SHOULD_RELEASE`) leaked into `QueueStatus.res`: "no error,
* but the packet should still be released". Firmware 2.8+ returns it for self-addressed packets, which are
* delivered through the synchronous local loopback instead of the TX queue — a success. Note it numerically
* collides with `Routing.Error.PKI_UNKNOWN_PUBKEY` (35); `QueueStatus.res` carries ErrorCode semantics, not
* Routing.Error.
*/
Tff7b72private Tff7b72const Tff7b72val Te6edf3ERRNO_SHOULD_RELEASE Tff7b72= T79c0ff3T79c0ff5
Tb4b4b4}
Tff7b72private Tff7b72var Te6edf3queueJobTb4b4b4: Te6edf3Job? Tff7b72= Tff7b72null
Tff7b72private Tff7b72val Te6edf3queueMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3queuedPackets Tff7b72= Te6edf3mutableListOfTff7b72<Te6edf3MeshPacketTff7b72>Tb4b4b4(Tb4b4b4)
T8b949e// Set to true by stopPacketQueue() under queueMutex. Checked by startPacketQueueLocked()
T8b949e// and the queue processor's finally block to prevent restarting a stopped queue.
Tff7b72private Tff7b72var Te6edf3queueStopped Tff7b72= Tff7b72false
Tff7b72private Tff7b72val Te6edf3responseMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3queueResponse Tff7b72= Te6edf3mutableMapOfTff7b72<Tffa657IntTb4b4b4, Te6edf3CompletableDeferredTff7b72<Tffa657BooleanTff7b72>Tff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3routingResponse Tff7b72= Te6edf3mutableMapOfTff7b72<Tffa657IntTb4b4b4, Te6edf3CompletableDeferredTff7b72<Tffa657BooleanTff7b72>Tff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3timeoutMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3sendAckTimeoutJobs Tff7b72= Te6edf3mutableMapOfTff7b72<Tffa657IntTb4b4b4, Te6edf3JobTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffsendToRadioTb4b4b4(Te6edf3pTb4b4b4: Te6edf3ToRadioTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffSending to radio Tffd700${Te6edf3pTb4b4b4.Te6edf3toPIIStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72val Te6edf3b Tff7b72= Te6edf3pTb4b4b4.Te6edf3encodeTb4b4b4(Tb4b4b4)
Te6edf3radioInterfaceServiceTb4b4b4.Te6edf3sendToRadioTb4b4b4(Te6edf3bTb4b4b4)
Te6edf3pTb4b4b4.Te6edf3packetTff7b72?.Te6edf3idTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3changeStatusTb4b4b4(Tffa657itTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ENROUTETb4b4b4) Tb4b4b4}
Tff7b72val Te6edf3packet Tff7b72= Te6edf3pTb4b4b4.Te6edf3packet
Tff7b72if Tb4b4b4(Te6edf3packetTff7b72?.Te6edf3decoded Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3packetToSave Tff7b72=
Te6edf3MeshLogTb4b4b4(
Te6edf3uuid Tff7b72= Te6edf3UuidTb4b4b4.Te6edf3randomTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3toStringTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3message_type Tff7b72= Ta5d6ff"Ta5d6ffPacketTa5d6ff"Tb4b4b4,
Te6edf3received_date Tff7b72= Te6edf3nowMillisTb4b4b4,
Te6edf3raw_message Tff7b72= Te6edf3packetTb4b4b4.Te6edf3toStringTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3fromNum Tff7b72= Te6edf3MeshLogTb4b4b4.Te6edf3NODE_NUM_LOCALTb4b4b4,
Te6edf3portNum Tff7b72= Te6edf3packetTb4b4b4.Te6edf3decodedTff7b72?.Te6edf3portnumTff7b72?.Te6edf3value Tff7b72?: T79c0ff0Tb4b4b4,
Te6edf3fromRadio Tff7b72= Te6edf3FromRadioTb4b4b4(Te6edf3packet Tff7b72= Te6edf3packetTb4b4b4)Tb4b4b4,
Tb4b4b4)
Te6edf3insertMeshLogTb4b4b4(Te6edf3packetToSaveTb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Enqueue [packet] for transmission. Order is preserved for sequential calls from the same coroutine (mutex
* acquisition is uncontested between sequential calls). Transactional sequences that require strict ordering across
* multiple calls — e.g. an `editSettings { … }` begin → writes → commit sequence — MUST be issued from a single
* coroutine; concurrent senders share FIFO only at the per-call grain.
*/
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffsendToRadioTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4) Tb4b4b4{
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueStopped Tff7b72= Tff7b72false
Te6edf3queuedPacketsTb4b4b4.Te6edf3addTb4b4b4(Te6edf3packetTb4b4b4)
Te6edf3startPacketQueueLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffSwallowedExceptionTa5d6ff"Tb4b4b4)
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffsendToRadioAndAwaitTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3connectionStateProviderTb4b4b4.Te6edf3connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadioAndAwait packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff skipped: not connectedTa5d6ff" Tb4b4b4}
Tff7b72return Tff7b72false
Tb4b4b4}
T8b949e// QueueStatus(res=0) only means that firmware queued the packet. Keep a separate
T8b949e// routing response so this strict caller waits for the later Routing ACK/NAK.
Tff7b72val Te6edf3deferred Tff7b72= Te6edf3CompletableDeferredTff7b72<Tffa657BooleanTff7b72>Tb4b4b4(Tb4b4b4)
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3routingResponseTff7b72[Te6edf3packetTb4b4b4.Te6edf3idTff7b72] Tff7b72= Te6edf3deferred Tb4b4b4}
Tff7b72return Tff7b72try Tb4b4b4{
Te6edf3sendToRadioTb4b4b4(Te6edf3packetTb4b4b4)
Te6edf3withTimeoutTb4b4b4(Te6edf3TIMEOUTTb4b4b4) Tb4b4b4{ Te6edf3deferredTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3TimeoutCancellationExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadioAndAwait packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff timeoutTa5d6ff" Tb4b4b4}
Tff7b72false
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e T8b949e// Preserve structured concurrency cancellation propagation.
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadioAndAwait packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff failed: Tffd700${Te6edf3eTb4b4b4.Te6edf3messageTffd700}Ta5d6ff" Tb4b4b4}
Tff7b72false
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3routingResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffstopPacketQueueTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// Run async so callers (non-suspend) don't block, but all mutations are
T8b949e// serialized under the same mutexes used by the queue processor and senders.
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffStopping packet queueJobTa5d6ff" Tb4b4b4}
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueStopped Tff7b72= Tff7b72true
Te6edf3queueJobTff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3queueJob Tff7b72= Tff7b72null
Te6edf3queuedPacketsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueResponseTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3forEach Tb4b4b4{ Tff7b72if Tb4b4b4(Tff7b72!Tffa657itTb4b4b4.Te6edf3isCompletedTb4b4b4) Tffa657itTb4b4b4.Te6edf3completeTb4b4b4(Tff7b72falseTb4b4b4) Tb4b4b4}
Te6edf3queueResponseTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3routingResponseTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3forEach Tb4b4b4{ Tff7b72if Tb4b4b4(Tff7b72!Tffa657itTb4b4b4.Te6edf3isCompletedTb4b4b4) Tffa657itTb4b4b4.Te6edf3completeTb4b4b4(Tff7b72falseTb4b4b4) Tb4b4b4}
Te6edf3routingResponseTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffhandleQueueStatusTb4b4b4(Te6edf3queueStatusTb4b4b4: Te6edf3QueueStatusTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ff[queueStatus] Tffd700${Te6edf3queueStatusTb4b4b4.Te6edf3toOneLineStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72val Tb4b4b4(Te6edf3successTb4b4b4, Te6edf3isFullTb4b4b4, Te6edf3requestIdTb4b4b4) Tff7b72=
Te6edf3withTb4b4b4(Te6edf3queueStatusTb4b4b4) Tb4b4b4{ Te6edf3TripleTb4b4b4(Te6edf3res Tff7b72=Tff7b72= T79c0ff0 Tff7b72|Tff7b72| Te6edf3res Tff7b72=Tff7b72= Te6edf3ERRNO_SHOULD_RELEASETb4b4b4, Te6edf3free Tff7b72=Tff7b72= T79c0ff0Tb4b4b4, Te6edf3mesh_packet_idTb4b4b4) Tb4b4b4}
T8b949e// Only the plain res==0 "queue accepted, now full" echo is skipped here. ERRNO_SHOULD_RELEASE denotes a
T8b949e// synchronous local-loopback delivery that still needs its queueResponse completed, even when free==0, or it
T8b949e// would hang until TIMEOUT.
Tff7b72if Tb4b4b4(Te6edf3queueStatusTb4b4b4.Te6edf3res Tff7b72=Tff7b72= T79c0ff0 Tff7b72&Tff7b72& Te6edf3isFullTb4b4b4) Tff7b72return
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3requestId Tff7b72!Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3requestIdTb4b4b4)Tff7b72?.Te6edf3completeTb4b4b4(Te6edf3successTb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3success Tff7b72|Tff7b72| Te6edf3queueStatusTb4b4b4.Te6edf3res Tff7b72=Tff7b72= Te6edf3ERRNO_SHOULD_RELEASETb4b4b4) Tb4b4b4{
Te6edf3routingResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3requestIdTb4b4b4)Tff7b72?.Te6edf3completeTb4b4b4(Te6edf3successTb4b4b4)
Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3queueResponseTb4b4b4.Te6edf3entries
Tb4b4b4.Te6edf3firstOrNull Tb4b4b4{ Tff7b72!Tffa657itTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3isCompleted Tb4b4b4}
Tff7b72?.Te6edf3let Tb4b4b4{ Tb4b4b4(Te6edf3packetIdTb4b4b4, Te6edf3responseTb4b4b4) Tff7b72-Tff7b72>
Te6edf3responseTb4b4b4.Te6edf3completeTb4b4b4(Te6edf3successTb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3success Tff7b72|Tff7b72| Te6edf3queueStatusTb4b4b4.Te6edf3res Tff7b72=Tff7b72= Te6edf3ERRNO_SHOULD_RELEASETb4b4b4) Tb4b4b4{
Te6edf3routingResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetIdTb4b4b4)Tff7b72?.Te6edf3completeTb4b4b4(Te6edf3successTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffremoveResponseTb4b4b4(Te6edf3dataRequestIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3completeTb4b4b4: Tffa657BooleanTb4b4b4) Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3dataRequestIdTb4b4b4)Tff7b72?.Te6edf3completeTb4b4b4(Te6edf3completeTb4b4b4)
Te6edf3routingResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3dataRequestIdTb4b4b4)Tff7b72?.Te6edf3completeTb4b4b4(Te6edf3completeTb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Starts the packet queue processor. Must be called while holding [queueMutex] to ensure the check-then-start is
* atomic — preventing two concurrent callers from launching duplicate processors.
*/
Tff7b72private Tff7b72fun Td2a8ffstartPacketQueueLockedTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3queueStoppedTb4b4b4) Tff7b72return
Tff7b72if Tb4b4b4(Te6edf3queueJobTff7b72?.Te6edf3isActive Tff7b72=Tff7b72= Tff7b72trueTb4b4b4) Tff7b72return
Te6edf3queueJob Tff7b72=
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Tff7b72try Tb4b4b4{
Tff7b72while Tb4b4b4(Te6edf3connectionStateProviderTb4b4b4.Te6edf3connectionStateTb4b4b4.Te6edf3value Tff7b72=Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3packet Tff7b72= Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queuedPacketsTb4b4b4.Te6edf3removeFirstOrNullTb4b4b4(Tb4b4b4) Tb4b4b4} Tff7b72?: Tff7b72break
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffSwallowedExceptionTa5d6ff"Tb4b4b4)
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3response Tff7b72= Te6edf3sendPacketTb4b4b4(Te6edf3packetTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff waitingTa5d6ff" Tb4b4b4}
Tff7b72val Te6edf3success Tff7b72= Te6edf3withTimeoutTb4b4b4(Te6edf3TIMEOUTTb4b4b4) Tb4b4b4{ Te6edf3responseTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff success Tffd700$Te6edf3successTa5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3TimeoutCancellationExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff timeoutTa5d6ff" Tb4b4b4}
T8b949e// Clean up the transport-queue deferred for this packet.
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e T8b949e// Preserve structured concurrency cancellation propagation.
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff failedTa5d6ff" Tb4b4b4}
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4) Tb4b4b4}
Tb4b4b4}
T8b949e// Deferred cleanup is now handled in the catch blocks above.
T8b949e// handleQueueStatus (normal success) and stopPacketQueue (bulk cleanup)
T8b949e// also remove entries, and these removals are idempotent.
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
T8b949e// Hold queueMutex so that clearing queueJob and the restart decision are
T8b949e// atomic with respect to new senders calling startPacketQueueLocked().
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueJob Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3queueStopped Tff7b72&Tff7b72& Te6edf3queuedPacketsTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3startPacketQueueLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffchangeStatusTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3mTb4b4b4: Te6edf3MessageStatusTb4b4b4) Tff7b72= Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3packetId Tff7b72!Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3getDataPacketByIdTb4b4b4(Te6edf3packetIdTb4b4b4)Tff7b72?.Te6edf3let Tb4b4b4{ Te6edf3p Tff7b72-Tff7b72>
Tff7b72if Tb4b4b4(Te6edf3pTb4b4b4.Te6edf3status Tff7b72!Tff7b72= Te6edf3mTb4b4b4) Tb4b4b4{
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3updateMessageStatusTb4b4b4(Te6edf3pTb4b4b4, Te6edf3mTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3m Tff7b72=Tff7b72= Te6edf3MessageStatusTb4b4b4.Te6edf3ENROUTETb4b4b4) Tb4b4b4{
Te6edf3scheduleSendAckTimeoutTb4b4b4(Te6edf3packetIdTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffrearmSendAckTimeoutsTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3getEnroutePacketsTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3p Tff7b72-Tff7b72>
Tff7b72val Te6edf3remaining Tff7b72= Te6edf3pTb4b4b4.Te6edf3time Tff7b72+ Te6edf3SEND_ACK_TIMEOUTTb4b4b4.Te6edf3inWholeMilliseconds Tff7b72- Te6edf3nowMillis
Te6edf3scheduleSendAckTimeoutTb4b4b4(Te6edf3pTb4b4b4.Te6edf3idTb4b4b4, Te6edf3remainingTb4b4b4.Te6edf3millisecondsTb4b4b4.Te6edf3coerceAtLeastTb4b4b4(Te6edf3REARM_GRACETb4b4b4)Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* A send whose routing ACK/NAK never reaches the app (typically because it was disconnected when the radio's
* response arrived) would stay ENROUTE — "Sending…" — forever. Stamp it as a retryable timeout instead; a late ACK
* still upgrades it via handleAckNak.
*
* One timer per packet: re-arming supersedes the pending one, so repeated reconnects cannot pile up timers for the
* same send. Timers deliberately survive a disconnect — the ack genuinely never arrived, and the resulting state is
* retryable — so a user who never reconnects still sees the send resolve.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffscheduleSendAckTimeoutTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3delayForTb4b4b4: Te6edf3Duration Tff7b72= Te6edf3SEND_ACK_TIMEOUTTb4b4b4) Tb4b4b4{
Te6edf3timeoutMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3sendAckTimeoutJobsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetIdTb4b4b4)Tff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3sendAckTimeoutJobsTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3removeAll Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3isCompleted Tb4b4b4}
Te6edf3sendAckTimeoutJobsTff7b72[Te6edf3packetIdTff7b72] Tff7b72=
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3delayTb4b4b4(Te6edf3delayForTb4b4b4)
T8b949e// Conditional in the DAO transaction: an ACK/NAK landing while this timer waited must win.
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3timeOutEnroutePacketTb4b4b4(Te6edf3packetIdTb4b4b4, Te6edf3RoutingTb4b4b4.Te6edf3ErrorTb4b4b4.Te6edf3TIMEOUTTb4b4b4.Te6edf3valueTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffgetDataPacketByIdTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4)Tb4b4b4: Te6edf3DataPacket? Tff7b72= Te6edf3withTimeoutOrNullTb4b4b4(T79c0ff1.Te6edf3secondsTb4b4b4) Tb4b4b4{
Tff7b72var Te6edf3dataPacketTb4b4b4: Te6edf3DataPacket? Tff7b72= Tff7b72null
Tff7b72while Tb4b4b4(Te6edf3dataPacket Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3dataPacket Tff7b72= Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3getPacketByIdTb4b4b4(Te6edf3packetIdTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3dataPacket Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Te6edf3delayTb4b4b4(T79c0ff1T79c0ff0T79c0ff0.Te6edf3millisecondsTb4b4b4)
Tb4b4b4}
Te6edf3dataPacket
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffsendPacketTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4)Tb4b4b4: Te6edf3DeferredTff7b72<Tffa657BooleanTff7b72> Tb4b4b4{
T8b949e// Register the transport-queue response before sending so an immediate QueueStatus cannot be missed.
Tff7b72val Te6edf3deferred Tff7b72= Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTb4b4b4.Te6edf3getOrPutTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4) Tb4b4b4{ Te6edf3CompletableDeferredTb4b4b4(Tb4b4b4) Tb4b4b4} Tb4b4b4}
Tff7b72try Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3connectionStateProviderTb4b4b4.Te6edf3connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3RadioNotConnectedExceptionTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3sendToRadioTb4b4b4(Te6edf3ToRadioTb4b4b4(Te6edf3packet Tff7b72= Te6edf3packetTb4b4b4)Tb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3exTb4b4b4: Te6edf3RadioNotConnectedExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3exTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio skipped: Not connected to radioTa5d6ff" Tb4b4b4}
Te6edf3removeResponseTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3complete Tff7b72= Tff7b72falseTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3exTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3exTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio error: Tffd700${Te6edf3exTb4b4b4.Te6edf3messageTffd700}Ta5d6ff" Tb4b4b4}
Te6edf3removeResponseTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4, Te6edf3complete Tff7b72= Tff7b72falseTb4b4b4)
Tb4b4b4}
T8b949e// Return a read-only Deferred view (kotlinx.coroutines 1.11+) so callers can await it
T8b949e// without being able to complete the underlying CompletableDeferred; cancellation is
T8b949e// still exposed via Deferred/Job.
Tff7b72return Te6edf3deferredTb4b4b4.Te6edf3asDeferredTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffinsertMeshLogTb4b4b4(Te6edf3packetToSaveTb4b4b4: Te6edf3MeshLogTb4b4b4) Tb4b4b4{
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{
Ta5d6ff"Ta5d6ffinsert: Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3message_typeTffd700}Ta5d6ff = Ta5d6ff" Tff7b72+
Ta5d6ff"Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3raw_messageTb4b4b4.Te6edf3toOneLineStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff from=Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3fromNumTffd700}Ta5d6ff"
Tb4b4b4}
Te6edf3meshLogRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3insertTb4b4b4(Te6edf3packetToSaveTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Served by rngit 1.5.0 - Generated in 0.09s